Skip to content

Fix race condition in DirectRunner executor shutdown - #39096

Merged
Amar3tto merged 2 commits into
apache:masterfrom
shunping:fix-flaky-rrio-test
Jun 25, 2026
Merged

Fix race condition in DirectRunner executor shutdown#39096
Amar3tto merged 2 commits into
apache:masterfrom
shunping:fix-flaky-rrio-test

Conversation

@shunping

Copy link
Copy Markdown
Collaborator

Previously, in ExecutorServiceParallelExecutor, if an exception occurred during registry cleanup (such as a timeout inside DoFn teardown), the pipeline state was transitioned to terminal before the exception was queued in visibleUpdates.

This introduced a race condition where waitUntilFinish() could detect the terminal state and exit successfully before the exception was offered to the updates queue, swallowing the exception.

This caused tests like CallTest.givenTeardownTimeout_throwsError to fail since they expected the pipeline to throw an exception. This change swaps the order so that the exception is posted to visibleUpdates before updating the pipeline state to terminal, ensuring the exception is always propagated.

fixes #38845

Previously, in ExecutorServiceParallelExecutor, if an exception
occurred during registry cleanup (such as a timeout inside
DoFn teardown), the pipeline state was transitioned to terminal
before the exception was queued in `visibleUpdates`.

This introduced a race condition where `waitUntilFinish()` could
detect the terminal state and exit successfully before the
exception was offered to the updates queue, swallowing the exception.

This caused tests like `CallTest.givenTeardownTimeout_throwsError`
to fail since they expected the pipeline to throw an exception.
This change swaps the order so that the exception is posted to
`visibleUpdates` before updating the pipeline state to terminal,
ensuring the exception is always propagated.
@shunping

Copy link
Copy Markdown
Collaborator Author

Run Java_IOs_Direct PreCommit

@shunping

Copy link
Copy Markdown
Collaborator Author

Running Java_IOs_Direct PreCommit with the change: https://github.com/apache/beam/actions/runs/28125966317

@shunping
shunping marked this pull request as ready for review June 24, 2026 20:03
@gemini-code-assist

Copy link
Copy Markdown
Contributor

Summary of Changes

Hello, I'm Gemini Code Assist1! I'm currently reviewing this pull request and will post my feedback shortly. In the meantime, here's a summary to help you and other reviewers quickly get up to speed!

This pull request addresses a race condition in the DirectRunner's executor shutdown logic. By ensuring that exceptions are properly queued before the pipeline state transitions to a terminal status, it guarantees that errors are correctly propagated to callers of waitUntilFinish, resolving intermittent test failures related to teardown timeouts.

Highlights

  • Race Condition Resolution: Reordered the pipeline state transition and exception propagation in the executor shutdown process to prevent exceptions from being swallowed.
  • Exception Propagation: Ensured that any errors occurring during registry cleanup are queued in visibleUpdates before the pipeline state is marked as terminal.
New Features

🧠 You can now enable Memory (public preview) to help Gemini Code Assist learn from your team's feedback. This makes future code reviews more consistent and personalized to your project's style. Click here to enable Memory in your admin console.

Using Gemini Code Assist

The full guide for Gemini Code Assist can be found on our documentation page, here are some quick tips.

Invoking Gemini

You can request assistance from Gemini at any point by creating a comment using either /gemini <command> or @gemini-code-assist <command>. Below is a summary of the supported commands on the current page.

Feature Command Description
Code Review /gemini review Performs a code review for the current pull request in its current state.
Pull Request Summary /gemini summary Provides a summary of the current pull request in its current state.
Comment @gemini-code-assist Responds in comments when explicitly tagged, both in pull request comments and review comments.
Help /gemini help Displays a list of available commands.

Customization

To customize the Gemini Code Assist for GitHub experience, repository maintainers can create a configuration file and/or provide a custom code review style guide (such as PEP-8 for Python) by creating and adding files to a .gemini/ folder in the base of the repository. Detailed instructions can be found here.

Limitations & Feedback

Gemini Code Assist may make mistakes. Please leave feedback on any instances where its feedback is incorrect or counterproductive. You can react with 👍 and 👎 on @gemini-code-assist comments. If you're interested in giving your feedback about your experience with Gemini Code Assist for GitHub and other Google products, sign up here.

Footnotes

  1. Review the Privacy Notices, Generative AI Prohibited Use Policy, Terms of Service, and learn how to configure Gemini Code Assist in GitHub here. Gemini can make mistakes, so double check it and use code with caution.

@gemini-code-assist gemini-code-assist Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Code Review

This pull request modifies the shutdown logic in ExecutorServiceParallelExecutor.java by deferring the pipeline state transition until after error aggregation and reporting. However, this introduces a critical issue: if an exception occurs during error formatting (such as a NullPointerException if an error message is null), the pipeline state transition will be bypassed, leading to an indefinite hang. To resolve this, the state transition should be executed within a finally block to guarantee it always runs.

Important

The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.

@shunping

Copy link
Copy Markdown
Collaborator Author

Run Java_IOs_Direct PreCommit

@shunping

shunping commented Jun 24, 2026

Copy link
Copy Markdown
Collaborator Author

Re-running the failed Java_IOs_Direct PreCommit workflow with the change: https://github.com/apache/beam/actions/runs/28131137789, and the run is successful.

@shunping

Copy link
Copy Markdown
Collaborator Author

r: @Amar3tto

@github-actions

Copy link
Copy Markdown
Contributor

Stopping reviewer notifications for this pull request: review requested by someone other than the bot, ceding control. If you'd like to restart, comment assign set of reviewers

@Amar3tto
Amar3tto merged commit 2352106 into apache:master Jun 25, 2026
21 checks passed
Amar3tto added a commit that referenced this pull request Jun 26, 2026
[release-2.75] Cherry-pick #39096 to the branch
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

The PreCommit Java IOs Direct job is flaky

3 participants